安装
安装包下载
链接:https://pan.baidu.com/s/1nBrOPbvxn-tVq1DBwEkkaw
提取码:7b1i
--来自百度网盘超级会员V5的分享
解压
tar -zxvf flink-1.16.0-bin-scala_2.12.tgz -C ../module/
提交作用的不同方式
Standalone运行模式
修改配置文件
vim flink-conf.yaml
JobManager节点地址(公共部分)
jobmanager.rpc.address: 10.240.8.68
jobmanager.bind-host: 0.0.0.0
rest.address: 10.240.8.68
rest.bind-address: 0.0.0.0
taskmanager.numberOfTaskSlots: 3
修改workers
vi workers
10.240.8.68
10.240.8.229
10.240.8.185
修改masters
vi masters
10.240.8.68:8081
同步安装包多所有的集群
./xsync.sh /home/bigdata/module/flink-1.16.0
TaskManager节点地址.需要配置为当前机器名(不同的work节点配置的host不一样就行了)
taskmanager.bind-host: 0.0.0.0
taskmanager.host: 10.240.8.229
启动集群
./start-cluster.sh
启动效果
master StandaloneSessionClusterEntrypoint
work TaskManagerRunner
Yarn运行模式
配置环境变量
那台机器要启动flink就配置那台就行了,不用所有的hadoop集群都要配置这个
sudo vim /etc/profile.d/my_env.sh
export HADOOP_HOME=/home/bigdata/module/hadoop-3.2.3
export PATH=$PATH:$HADOOP_HOME/bin
export PATH=$PATH:$HADOOP_HOME/sbin
export HADOOP_CONF_DIR=${HADOOP_HOME}/etc/hadoop
export HADOOP_CLASSPATH=`hadoop classpath`
source /etc/profile.d/my_env.sh
YARN Session模式
执行脚本命令向YARN集群申请资源,开启一个YARN会话,启动Flink集群。
bin/yarn-session.sh -nm test
- -d:分离模式,如果你不想让Flink YARN客户端一直前台运行,可以使用这个参数,即使关掉当前对话窗口,YARN session也可以后台运行。
- -jm(--jobManagerMemory):配置JobManager所需内存,默认单位MB。
- -nm(--name):配置在YARN UI界面上显示的任务名。
- -qu(--queue):指定YARN队列名。
- -tm(--taskManager):配置每个TaskManager所使用内存。
YARN的会话模式不会把集群资源固定,同样是动态分配的。
向集群提交任务
bin/flink run -c com.xx.xx xxx.jar
比如官网测试程序
./bin/flink run ./examples/streaming/TopSpeedWindowing.jar
YARN Per-job模式
在YARN环境中,由于有了外部平台做资源调度,所以我们也可以直接向YARN提交一个单独的作业,从而启动一个Flink集群。
./bin/flink run -t yarn-per-job --detached xx.jar
比如官网测试用例
./bin/flink run -t yarn-per-job --detached ./examples/streaming/TopSpeedWindowing.jar
yarn application -kill application_1698659960145_0005
yarn logs -applicationId application_1698659960145_0005
注意:如果启动过程中报如下异常。
Exception in thread “Thread-5” java.lang.IllegalStateException: Trying to access closed classloader. Please check if you store classloaders directly or indirectly in static fields. If the stacktrace suggests that the leak occurs in a third party library and cannot be fixed immediately, you can disable this check with the configuration ‘classloader.check-leaked-classloader’.
at org.apache.flink.runtime.execution.librarycache.FlinkUserCodeClassLoaders
解决方法:
vim flink-conf.yaml
classloader.check-leaked-classloader: false
YARN Application模式
可以通过yarn.provided.lib.dirs配置选项指定位置,将flink的依赖上传到远程。
- 上传flink的lib和plugins到HDFS上。
hadoop fs -mkdir /flink-dist
hadoop fs -put /home/bigdata/module/flink-1.16.0/lib/ /flink-dist
hadoop fs -put /home/bigdata/module/flink-1.16.0/plugins/ /flink-dist
- 上传自己的jar包到HDFS。
hadoop fs -mkdir /flink-jars
hadoop fs -put ./examples/streaming/TopSpeedWindowing.jar /flink-jars
- 提交作业
bin/flink run-application -t yarn-application -Dyarn.provided.lib.dirs="hdfs://bigdatacluster/flink-dist" -c com.xxx hdfs://bigdatacluster/flink-jars/xx.jar
比如
bin/flink run-application -t yarn-application -Dyarn.provided.lib.dirs="hdfs://bigdatacluster/flink-dist" hdfs://bigdatacluster/flink-jars/TopSpeedWindowing.jar
这种方式下,flink本身的依赖和用户jar可以预先上传到HDFS,而不需要单独发送到集群,这就使得作业提交更加轻量了。
配置历史服务器
运行 Flink job 的集群一旦停止,只能去 yarn 或本地磁盘上查看日志,不再可以查看作业挂掉之前的运行的 Web UI,很难清楚知道作业在挂的那一刻到底发生了什么。
创建对应的目录
hadoop fs -mkdir -p /logs/flink-job
修改配置文件
vi /home/bigdata/module/flink-1.16.0/conf/flink-conf.yaml
jobmanager.archive.fs.dir: hdfs://bigdatacluster/logs/flink-job
historyserver.web.address: 0.0.0.0
historyserver.web.port: 8082
historyserver.archive.fs.dir: hdfs://bigdatacluster/logs/flink-job
historyserver.archive.fs.refresh-interval: 5000
启停历史服务器
bin/historyserver.sh start
提交一个作业
bin/flink run-application -t yarn-application -Dyarn.provided.lib.dirs="hdfs://bigdatacluster/flink-dist" hdfs://bigdatacluster/flink-jars/TopSpeedWindowing.jar
访问
http://master1:8082
停止历史服务器
bin/historyserver.sh stop